Skip to content

[REA-5060] Add the managed distributed video model adapter - #176

Closed
maxmb1 wants to merge 1 commit into
codex/rea-5060-worker-groupfrom
codex/rea-5060-distributed-video-model
Closed

maxmb1 wants to merge 1 commit into
codex/rea-5060-worker-groupfrom
codex/rea-5060-distributed-video-model

Conversation

@maxmb1

@maxmb1 maxmb1 commented Sep 9, 2026

Copy link
Copy Markdown

Why

Model authors should declare their worker and typed video output without reimplementing threading, session bookkeeping, and cleanup around WorkerGroup. This layer supplies that managed authoring path without introducing an engine-host dependency.

What Changed

Adds the experimental DistributedVideoModel adapter with worker configuration, control snapshots, frame-to-output mapping, pause/restart, and worker cleanup. Authoring and configuration tests cover the adapter contract, while the dependent example PR exercises the complete run-loop lifecycle. The scope remains single-node, single-session uint8 video, not arbitrary structured outputs or multi-session scheduling. Local Python 3.12 unit/contract tests, lint, and type checking pass.

## Why

Model authors should declare their worker and typed video output without reimplementing threading, session bookkeeping, and cleanup around WorkerGroup. This layer supplies that managed authoring path without introducing an engine-host dependency.

## What Changed

Adds the experimental DistributedVideoModel adapter with worker configuration, control snapshots, frame-to-output mapping, pause/restart, and worker cleanup. Authoring and configuration tests cover the adapter contract, while the dependent example PR exercises the complete run-loop lifecycle. The scope remains single-node, single-session uint8 video, not arbitrary structured outputs or multi-session scheduling. Local Python 3.12 unit/contract tests, lint, and type checking pass.

Signed-off-by: Kavya <58096384+maxmb1@users.noreply.github.com>

maxmb1 commented Sep 9, 2026

Copy link
Copy Markdown
Author

Warning

This pull request is not mergeable via GitHub because a downstack PR is open. Once all requirements are satisfied, merge this PR as a stack on Graphite.
Learn more

This stack of pull requests is managed by Graphite. Learn more about stacking.

@Dere-Wah

Dere-Wah commented Sep 17, 2026 •

Copy link
Copy Markdown
Contributor

Closing this in favour of a new stack that lands the same tool with a smaller contract. Design: DistributedWorkers; tracking: REA-6343 (world_size=1) and REA-6345 (world_size > 1).

Most of what is here carries over as is: the shared-memory block and pin(), liveness polling inside every wait, the _healthy flag with terminate-then-kill teardown, the child environment and faulthandler, per-rank log demotion, the world-size check in the child, and the picklability check on what crosses spawn. That work is the base of the new stack.

What changes: the worker's methods become the model half's three names (load / generate / reset), the runner accepts any class with them, the thing that crosses is the author's own struct carried through pickle protocol-5 out-of-band buffers instead of a fixed frame buffer, every rank's answer counts before the outcome is decided, and one worker still runs in its own process. DistributedVideoModel is dropped because the runtime's step loop is that layer now, and world_size is a constructor argument rather than a manifest read.

New PRs are linked from REA-6343 and REA-6345 as they open. Thanks for the groundwork; this stack stands on it.

The new stack: #195 (transport and wire), #196 (DistributedRunner at world_size=1), #197 (world_size > 1).

@Dere-Wah Dere-Wah closed this Sep 17, 2026
Dere-Wah added a commit that referenced this pull request Sep 24, 2026
…#195)

https://linear.app/reactor-team/issue/REA-6343 · Design: [DistributedWorkers](https://app.notion.com/p/reactor-technologies/DistributedWorkers-3decacd6a3968030b598d95d6df32ba0)

First of three. Standalone: nothing else in the runtime imports `reactor_runtime.distributed`, and this stack does not depend on the `ReactorApp` stack (#187 to #192). Lifted from #134, #175, #176, #177, which it closes.

## Why

The runner in #196 moves a model's input and result between processes. Both are the author's own objects and hold large arrays. Pickling them through a queue copies every byte; a fixed frame buffer forces a shape on the author. Neither is acceptable.

## What changed

`pack()` pickles with protocol 5. Each flat buffer (a contiguous numpy array) goes into a shared-memory block; the pickle stays small and rides the queue as a header that names the block and the spans. `unpack()` attaches to whatever block the header names and copies the spans out.

```python
slot = SharedSlot()                  # writer side, one per direction
header = pack(result, slot)          # ~1 KB on the queue; the arrays are in the block

reader = SlotReader()                # reader side
result = unpack(header, reader)      # the same object, arrays copied out
```

The block grows to the writer's high-water mark. Since the header names the block, growth needs no protocol. A torch tensor or a sliced array does not expose a flat buffer and pickles inline: slower, still correct, reported once per process. `SharedSlotAllocationFailed` names `--shm-size` and the pod's `/dev/shm` when a block cannot be created.

Also here: the wire (`Generate`, `Reset`, `Shutdown`, `Loaded`, `Answer`) and the runner's errors (`WorkerCrashed`, `WorkerTimeout`, `RankDesync`).

## Verification

`tests/unit/distributed/test_distributed_ipc.py`: a struct with a dict, str, enum, list, bytes, nested dataclass, `None`, and two arrays of different dtypes crosses intact with both arrays out of band; growth mid-call keeps what was already written and unlinks the old block; a non-contiguous array pickles inline and is reported once; the allocation error names the fix. `mise run lint`, `typecheck`, `test` green.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants